Micron Document
🎖️GitЯра🎖️


Displaying Raw • Download

core/takserver/src/commonMain/kotlin/org/meshtastic/core/takserver/TAKServerManager.kt renovate/net.java.dev.jna-jna-5.x (5a8f4966) Text, 7.09 KB

T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.takserver

Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableSharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableStateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.SharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.StateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asSharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asStateFlow
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlin.time.Clock
Tff7b72import T7ee787kotlin.time.Duration.Companion.minutes

T8b949e/** A CoT message received from a connected TAK client, paired with the client's identity. */
Tff7b72data Tff7b72class T56d364InboundCoTMessageTb4b4b4(Tff7b72val Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4, Tff7b72val Te6edf3clientInfoTb4b4b4: Te6edf3TAKClientInfo? Tff7b72= Tff7b72nullTb4b4b4)

Tff7b72interface T56d364TAKServerManager Tb4b4b4{
Tff7b72val Te6edf3isRunningTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72>
Tff7b72val Te6edf3connectionCountTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657IntTff7b72>
Tff7b72val Te6edf3inboundMessagesTb4b4b4: Te6edf3SharedFlowTff7b72<Te6edf3InboundCoTMessageTff7b72>

T8b949e/** Start the TAK server using [scope]. Port is fixed at [TAKServer] construction time. */
Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4)

Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4)

Tff7b72fun Td2a8ffbroadcastTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4)

T8b949e/** Broadcast raw XML verbatim to TAK clients, bypassing CoTMessage parsing. */
Tff7b72fun Td2a8ffbroadcastRawXmlTb4b4b4(Te6edf3xmlTb4b4b4: Tffa657StringTb4b4b4)
Tb4b4b4}

Tff7b72internal Tff7b72class T56d364TAKServerManagerImplTb4b4b4(Tff7b72private Tff7b72val Te6edf3takServerTb4b4b4: Te6edf3TAKServerTb4b4b4) Tb4b4b4: Te6edf3TAKServerManager Tb4b4b4{

Tff7b72private Tff7b72var Te6edf3scopeTb4b4b4: Te6edf3CoroutineScope? Tff7b72= Tff7b72null

Tff7b72private Tff7b72val Te6edf3_isRunning Tff7b72= Te6edf3MutableStateFlowTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72override Tff7b72val Te6edf3isRunningTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72> Tff7b72= Te6edf3_isRunningTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)

T8b949e// Mirror TAKServer's event-driven connection count — no polling needed
Tff7b72override Tff7b72val Te6edf3connectionCountTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657IntTff7b72> Tff7b72= Te6edf3takServerTb4b4b4.Te6edf3connectionCount

Tff7b72private Tff7b72val Te6edf3_inboundMessages Tff7b72= Te6edf3MutableSharedFlowTff7b72<Te6edf3InboundCoTMessageTff7b72>Tb4b4b4(Te6edf3extraBufferCapacity Tff7b72= T79c0ff6T79c0ff4Tb4b4b4)
Tff7b72override Tff7b72val Te6edf3inboundMessagesTb4b4b4: Te6edf3SharedFlowTff7b72<Te6edf3InboundCoTMessageTff7b72> Tff7b72= Te6edf3_inboundMessagesTb4b4b4.Te6edf3asSharedFlowTb4b4b4(Tb4b4b4)

T8b949e// Offline message queue — buffers mesh-originated CoT messages when no TAK
T8b949e// clients are connected, then drains them when a client reconnects. Entries
T8b949e// expire after OFFLINE_QUEUE_TTL to avoid delivering stale situational data.
Tff7b72private Tff7b72data Tff7b72class T56d364QueuedMessageTb4b4b4(Tff7b72val Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4, Tff7b72val Te6edf3enqueuedAtTb4b4b4: Te6edf3kotlinTb4b4b4.Te6edf3timeTb4b4b4.Te6edf3InstantTb4b4b4)

Tff7b72private Tff7b72val Te6edf3offlineQueue Tff7b72= Te6edf3ArrayDequeTff7b72<Te6edf3QueuedMessageTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3offlineQueueMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)

Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72private Tff7b72val Te6edf3OFFLINE_QUEUE_TTL Tff7b72= T79c0ff5.Te6edf3minutes
Tff7b72private Tff7b72const Tff7b72val Te6edf3OFFLINE_QUEUE_MAX_SIZE Tff7b72= T79c0ff5T79c0ff0
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTAKServerManager already runningTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
T8b949e// Assign scope AFTER the guard so a second concurrent start() can never
T8b949e// overwrite the active scope without actually restarting the server.
Tff7b72thisTb4b4b4.Te6edf3scope Tff7b72= Te6edf3scope

Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
T8b949e// Wire up inbound message handler BEFORE starting so no messages are lost.
T8b949e// Use tryEmit (non-suspending) with extraBufferCapacity to avoid launching a
T8b949e// new coroutine per message, which would create unbounded coroutines under
T8b949e// high message rates and could reorder messages.
Te6edf3takServerTb4b4b4.Te6edf3onMessage Tff7b72= Tb4b4b4{ Te6edf3cotMessageTb4b4b4, Te6edf3clientInfo Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_inboundMessagesTb4b4b4.Te6edf3tryEmitTb4b4b4(Te6edf3InboundCoTMessageTb4b4b4(Te6edf3cotMessageTb4b4b4, Te6edf3clientInfoTb4b4b4)Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK inbound message buffer full; dropping message from Tffd700${Te6edf3clientInfoTff7b72?.Te6edf3idTffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Te6edf3takServerTb4b4b4.Te6edf3onClientConnected Tff7b72= Tb4b4b4{ Te6edf3drainOfflineQueueTb4b4b4(Tb4b4b4) Tb4b4b4}

Tff7b72val Te6edf3result Tff7b72= Te6edf3takServerTb4b4b4.Te6edf3startTb4b4b4(Te6edf3scopeTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3resultTb4b4b4.Te6edf3isSuccessTb4b4b4) Tb4b4b4{
Te6edf3_isRunningTb4b4b4.Te6edf3value Tff7b72= Tff7b72true
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Server startedTa5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3resultTb4b4b4.Te6edf3exceptionOrNullTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to start TAK ServerTa5d6ff" Tb4b4b4}
T8b949e// Clear onMessage if start failed so we don't hold a reference unnecessarily
Te6edf3takServerTb4b4b4.Te6edf3onMessage Tff7b72= Tff7b72null
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Flip the running flag and null out the scope BEFORE stopping the server so
T8b949e// any broadcast()/drainOfflineQueue() that races stop() sees _isRunning=false
T8b949e// and exits early instead of launching coroutines on a scope that is about to
T8b949e// be discarded.
Te6edf3_isRunningTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Te6edf3scope Tff7b72= Tff7b72null
Te6edf3takServerTb4b4b4.Te6edf3onMessage Tff7b72= Tff7b72null
Te6edf3takServerTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Server stoppedTa5d6ff" Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffbroadcastTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3scopeTff7b72?.Te6edf3launch Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3takServerTb4b4b4.Te6edf3hasConnectionsTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3takServerTb4b4b4.Te6edf3broadcastTb4b4b4(Te6edf3cotMessageTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
T8b949e// No TAK clients connected — queue for delivery when one reconnects
Te6edf3offlineQueueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Evict expired entries
Tff7b72val Te6edf3cutoff Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4) Tff7b72- Te6edf3OFFLINE_QUEUE_TTL
Tff7b72while Tb4b4b4(Te6edf3offlineQueueTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4) Tff7b72&Tff7b72& Te6edf3offlineQueueTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3enqueuedAt Tff7b72< Te6edf3cutoffTb4b4b4) Tb4b4b4{
Te6edf3offlineQueueTb4b4b4.Te6edf3removeFirstTb4b4b4(Tb4b4b4)
Tb4b4b4}
T8b949e// Cap size to prevent unbounded growth
Tff7b72if Tb4b4b4(Te6edf3offlineQueueTb4b4b4.Te6edf3size Tff7b72>Tff7b72= Te6edf3OFFLINE_QUEUE_MAX_SIZETb4b4b4) Tb4b4b4{
Te6edf3offlineQueueTb4b4b4.Te6edf3removeFirstTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3offlineQueueTb4b4b4.Te6edf3addLastTb4b4b4(Te6edf3QueuedMessageTb4b4b4(Te6edf3cotMessageTb4b4b4, Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4)Tb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Tff7b72override Tff7b72fun Td2a8ffbroadcastRawXmlTb4b4b4(Te6edf3xmlTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3scopeTff7b72?.Te6edf3launch Tb4b4b4{ Te6edf3takServerTb4b4b4.Te6edf3broadcastRawXmlTb4b4b4(Te6edf3xmlTb4b4b4) Tb4b4b4}
Tb4b4b4}

T8b949e/**
* Drain any queued messages to the newly connected TAK client. Called by the server when a TAK client connects
* (Connected event).
*/
Tff7b72internal Tff7b72fun Td2a8ffdrainOfflineQueueTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3scopeTff7b72?.Te6edf3launch Tb4b4b4{
Tff7b72val Te6edf3messages Tff7b72=
Te6edf3offlineQueueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3cutoff Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4) Tff7b72- Te6edf3OFFLINE_QUEUE_TTL
Tff7b72val Te6edf3valid Tff7b72= Te6edf3offlineQueueTb4b4b4.Te6edf3filter Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3enqueuedAt Tff7b72>Tff7b72= Te6edf3cutoff Tb4b4b4}Tb4b4b4.Te6edf3map Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3cotMessage Tb4b4b4}
Te6edf3offlineQueueTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3valid
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3messagesTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffDraining Tffd700${Te6edf3messagesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff queued message(s) to reconnected TAK clientTa5d6ff" Tb4b4b4}
Te6edf3messagesTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3takServerTb4b4b4.Te6edf3broadcastTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}

Served by rngit 1.5.0 - Generated in 0.13s